KAFKA-19662: Reset share group offsets for unsubscribed topics#20708
Conversation
JimmyWang6
left a comment
There was a problem hiding this comment.
LGTM, thanks for the change! Just two small things left.
| ) { | ||
| List<CoordinatorRecord> records = new ArrayList<>(); | ||
| ShareGroup group = groupMetadataManager.shareGroup(groupId); | ||
| group.validateOffsetsAlterable(); |
There was a problem hiding this comment.
I think the method validateOffsetsAlterable() could be removed.
| } catch (ExecutionException ee) { | ||
| Throwable cause = ee.getCause(); | ||
| if (cause instanceof GroupIdNotFoundException) { | ||
| resetOffsetsForInactiveGroup(groupId); |
There was a problem hiding this comment.
Just wondering, why is there a retry mechanism when GroupIdNotFoundException is thrown? Is this for better reliability? I see ConsumerGroupCommand also does this.
There was a problem hiding this comment.
I think this PR aims to align the behavior between consumer groups and shared consumer groups.
For regular consumer groups, we already have similar logic — if the group doesn’t exist, it will be automatically created. see
In this PR, the shared consumer group now follows the same behavior: https://github.com/apache/kafka/pull/20708/files#diff-00f0f81cf13e66781777d94f7d2e68a581663385c37e98792507f2294c91bb09R8309-R8310
Back to your question, I think your understanding is correct. This is something like a retry mechanism: if describeShareGroups throws a GroupIdNotFoundException, then will call resetOffsetsForInactiveGroup again, which will implicitly create the shared group.
Perhaps Andrew can elaborate more 😃
There was a problem hiding this comment.
It's because if you do the operation on a non-existent group, it creates the group (which is by definition empty) and then resets the offsets. I know this is a bit weird, but if you want to initialise the group in a particular state before starting any consumers, this is how you would do it.
| return partitionsToReset; | ||
| } | ||
|
|
||
| private Map<TopicPartition, OffsetAndMetadata> prepareOffsetsToReset(String groupId, Collection<TopicPartition> partitionsToReset) { |
There was a problem hiding this comment.
nit: now groupId become unused, could you please remove this variable? thanks!
Yes, I agree that the test failure is related. I'm going to fix that test and probably add another few. The fact that so little failed is a good sign that overall this way of addressing the problem works well. |
apoorvmittal10
left a comment
There was a problem hiding this comment.
Thanks for the PR, LGTM!
…e#20708) This PR allows the kafka-share-groups.sh --reset-offsets tool to be used to set offsets for topics which are not currently subscribed in a share group. It also works if the share group does not yet exist. This brings the capability in line with the equivalent function in Kafka-consumer-groups.sh. The primary purpose is to allow offsets to be set before the share group is first used as a way of initialising in a known state. Reviewers: Jimmy Wang <48462172+JimmyWang6@users.noreply.github.com>, Kuan-Po Tseng <brandboat@gmail.com>, Apoorv Mittal <apoorvmittal10@gmail.com>
…e#20708) This PR allows the kafka-share-groups.sh --reset-offsets tool to be used to set offsets for topics which are not currently subscribed in a share group. It also works if the share group does not yet exist. This brings the capability in line with the equivalent function in Kafka-consumer-groups.sh. The primary purpose is to allow offsets to be set before the share group is first used as a way of initialising in a known state. Reviewers: Jimmy Wang <48462172+JimmyWang6@users.noreply.github.com>, Kuan-Po Tseng <brandboat@gmail.com>, Apoorv Mittal <apoorvmittal10@gmail.com>
…e#20708) This PR allows the kafka-share-groups.sh --reset-offsets tool to be used to set offsets for topics which are not currently subscribed in a share group. It also works if the share group does not yet exist. This brings the capability in line with the equivalent function in Kafka-consumer-groups.sh. The primary purpose is to allow offsets to be set before the share group is first used as a way of initialising in a known state. Reviewers: Jimmy Wang <48462172+JimmyWang6@users.noreply.github.com>, Kuan-Po Tseng <brandboat@gmail.com>, Apoorv Mittal <apoorvmittal10@gmail.com>
…e#20708) This PR allows the kafka-share-groups.sh --reset-offsets tool to be used to set offsets for topics which are not currently subscribed in a share group. It also works if the share group does not yet exist. This brings the capability in line with the equivalent function in Kafka-consumer-groups.sh. The primary purpose is to allow offsets to be set before the share group is first used as a way of initialising in a known state. Reviewers: Jimmy Wang <48462172+JimmyWang6@users.noreply.github.com>, Kuan-Po Tseng <brandboat@gmail.com>, Apoorv Mittal <apoorvmittal10@gmail.com>
…e#20708) This PR allows the kafka-share-groups.sh --reset-offsets tool to be used to set offsets for topics which are not currently subscribed in a share group. It also works if the share group does not yet exist. This brings the capability in line with the equivalent function in Kafka-consumer-groups.sh. The primary purpose is to allow offsets to be set before the share group is first used as a way of initialising in a known state. Reviewers: Jimmy Wang <48462172+JimmyWang6@users.noreply.github.com>, Kuan-Po Tseng <brandboat@gmail.com>, Apoorv Mittal <apoorvmittal10@gmail.com>
* Revert "Merge pull request #161 from showuon/cl3_ut" This reverts commit 55efc8d, reversing changes made to 1c8b6d3. * Revert "Merge pull request #155 from showuon/cl3_share_auto_create" This reverts commit 05f8874, reversing changes made to 143e6cc. * MINOR: Cleanups in Tools Module (3/n) (apache#20332) This PR aims at cleaning up the tools module further by getting rid of some extra code which can be replaced by `record` Reviewers: Chia-Ping Tsai <chia7712@gmail.com> * MINOR: kafka-stream-groups.sh should fail quickly if the partition leader is unavailable (apache#20271) This PR applies the same partition leader check for `StreamsGroupCommand` as `ShareGroupCommand` and `ConsumerGroupCommand` to avoid the command execution timeout. Reviewers: Lucas Brutschy <lucasbru@apache.org> * KAFKA-19662: Allow resetting offset for unsubscribed topic in kafka-share-groups.sh (apache#20453) The `kafka-share-groups.sh` tool checks whether a topic already has a start-offset in the share group when resetting offsets. This is not necessary. By removing the check, it is possible to set a start offset for a topic which has not yet but will be subscribed in the future, thus initialising the consumption point. There is still a small piece of outstanding work to do with resetting the offset for a non-existent group which should also create the group. A subsequent PR will be used to address that. Reviewers: Jimmy Wang <48462172+JimmyWang6@users.noreply.github.com>, Lan Ding <isDing_L@163.com>, Apoorv Mittal <apoorvmittal10@gmail.com> * KAFKA-19662: Reset share group offsets for unsubscribed topics (apache#20708) This PR allows the kafka-share-groups.sh --reset-offsets tool to be used to set offsets for topics which are not currently subscribed in a share group. It also works if the share group does not yet exist. This brings the capability in line with the equivalent function in Kafka-consumer-groups.sh. The primary purpose is to allow offsets to be set before the share group is first used as a way of initialising in a known state. Reviewers: Jimmy Wang <48462172+JimmyWang6@users.noreply.github.com>, Kuan-Po Tseng <brandboat@gmail.com>, Apoorv Mittal <apoorvmittal10@gmail.com> --------- Co-authored-by: Sanskar Jhajharia <sjhajharia@confluent.io> Co-authored-by: jimmy <wangzhiwang611@gmail.com> Co-authored-by: Andrew Schofield <aschofield@confluent.io>
This PR allows the kafka-share-groups.sh --reset-offsets tool to be used
to set offsets for topics which are not currently subscribed in a share
group. It also works if the share group does not yet exist. This brings
the capability in line with the equivalent function in
Kafka-consumer-groups.sh. The primary purpose is to allow offsets to be
set before the share group is first used as a way of initialising in a
known state.
Reviewers: Jimmy Wang 48462172+JimmyWang6@users.noreply.github.com,
Kuan-Po Tseng brandboat@gmail.com, Apoorv Mittal
apoorvmittal10@gmail.com